welo git mirrors
Commit 5c7626aabeda1edf540b71eb9112d13b569414e7
Parents : 4d1e6f1
Author : welo | main nixos computer <empty@empty.com>
Date : 2026-07-08T19:06:11+01:00
added awaits everywhere
Changes
5 files changed, 53 insertions(+), 40 deletions(-)
Diff
diff --git a/src/client.rs b/src/client.rs
index 62ee452..1b87cd5 100644
--- a/src/client.rs
+++ b/src/client.rs
@@ -50,14 +50,17 @@ pub async fn run_client(server_hex: &str, listen_addr: &str) {
pub async fn run_client_forward(server_hex: &str, ports: Vec<ForwardedPort> ) {
if let Some(destination) = decode_hash(server_hex).await {
if let Some((mux, node, rx)) = connect_rns(destination).await {
+ info!("a");
let reconnect_notify = reconnect_generator(mux.clone(), node, rx, destination).await;
+ info!("b");
for port in ports {
match port.r#type {
PortType::Tcp => {
tcp_tunnel(mux.clone(), reconnect_notify.clone(), port).await
},
PortType::Udp => {
+ info!("c");
udp_tunnel(mux.clone(), reconnect_notify.clone(), port).await
},
PortType::TcpUdp => {
@@ -218,7 +221,7 @@ async fn establish_link(
} => {
if lid == link_id && is_initiator {
info!("RNS link established (rtt={:.1}ms)", rtt * 1000.0);
- mux.set_link_id(link_id);
+ mux.set_link_id(link_id).await;
return true;
}
}
@@ -263,12 +266,12 @@ async fn dispatch_and_reconnect(
// println!("{:?}",data);
for frame in mux.receive_data(&data).await {
// println!("type {:?} sid: {:?}", frame.frame_type, frame.session_id);
- mux.dispatch(frame);
+ mux.dispatch(frame).await;
}
}
ProxyEvent::LinkClosed { link_id, reason } => {
warn!("Connection lost (link={}, reason={:?})", link_id, reason);
- mux.reset();
+ mux.reset().await;
break; // Exit dispatch loop to reconnect
}
_ => {}
@@ -369,7 +372,7 @@ pub async fn connect_tcp_server_side(
info!("[{}] -> {}:{} tcp", sid, host, port);
// Send CONNECT frame through RNS
let connect_payload = encode_connect_payload(&host, port,false);
- mux.send(FrameType::Connect, sid, connect_payload);
+ mux.send(FrameType::Connect, sid, connect_payload).await;
// Wait for CONN_OK or CONN_ERR with timeout
tokio::time::timeout(Duration::from_secs(15), async {
@@ -423,8 +426,8 @@ async fn handle_tcp_connect(
Ok(s) => s,
Err(e) => {
debug!("[{}] Failed to send SOCKS5 reply: {}", sid, e);
- mux.send(FrameType::Close, sid, Vec::new());
- mux.drop_session(sid);
+ mux.send(FrameType::Close, sid, Vec::new()).await;
+ mux.drop_session(sid).await;
return None;
}
}
@@ -432,13 +435,13 @@ async fn handle_tcp_connect(
Some(Err(reason)) => {
warn!("[{}] Remote connect failed: {}", sid, reason);
let _ = proto.reply_error(&ReplyError::GeneralFailure).await;
- mux.drop_session(sid);
+ mux.drop_session(sid).await;
return None;
}
None => {
warn!("[{}] Connect timeout", sid);
let _ = proto.reply_error(&ReplyError::TtlExpired).await;
- mux.drop_session(sid);
+ mux.drop_session(sid).await;
return None;
}
};
@@ -469,7 +472,7 @@ pub async fn udp_bind_connect(
// Send CONNECT frame through RNS
let connect_payload = encode_connect_payload(&host, port, true);
- mux.send(FrameType::Connect, sid, connect_payload);
+ mux.send(FrameType::Connect, sid, connect_payload).await;
// Wait for CONN_OK or CONN_ERR with timeout
tokio::time::timeout(Duration::from_secs(15), async {
@@ -515,8 +518,8 @@ async fn handle_udp_connect(
Ok(s) => Some((udp_stream,s)),
Err(e) => {
debug!("[{}] Failed to send SOCKS5 reply: {:?}", sid, e);
- mux.send(FrameType::Close, sid, Vec::new());
- mux.drop_session(sid);
+ mux.send(FrameType::Close, sid, Vec::new()).await;
+ mux.drop_session(sid).await;
return None;
}
}
@@ -524,13 +527,13 @@ async fn handle_udp_connect(
Ok(Err(reason)) => {
warn!("[{}] Remote connect failed: {}", sid, reason);
let _ = proto.reply_error(&ReplyError::GeneralFailure).await;
- mux.drop_session(sid);
+ mux.drop_session(sid).await;
return None;
}
Err(_) => {
warn!("[{}] Connect timeout", sid);
let _ = proto.reply_error(&ReplyError::TtlExpired).await;
- mux.drop_session(sid);
+ mux.drop_session(sid).await;
return None;
}
}
diff --git a/src/forwarding.rs b/src/forwarding.rs
index a95ff7b..90becc6 100644
--- a/src/forwarding.rs
+++ b/src/forwarding.rs
@@ -114,6 +114,7 @@ pub async fn udp_tunnel(mux: MuxHandle, reconnect_notify: Arc<Notify> , port: Fo
}
};
+ info!("d");
let reference = &mux;
@@ -123,23 +124,32 @@ pub async fn udp_tunnel(mux: MuxHandle, reconnect_notify: Arc<Notify> , port: Fo
let (stream, addr) = match accept_result {
Ok(sa) => sa,
Err(e) => {
+ info!("hi");
continue;
}
};
+ info!("5");
let mux = reference.clone();
+ info!("1");
if !mux.is_connected().await {
drop(stream);
continue;
}
+ info!("2");
let sid = mux.next_session_id().await;
+ info!("2");
let mut session_rx = mux.register_session(sid).await;
+ info!("2");
let mux_clone = mux.clone();
+ info!("2");
+ info!("2");
// let connect_result = udp_bind_connect(sid,mux.clone(), &mut session_rx, target_addr).await;
let target_addr = TargetAddr::Ip(SocketAddr::new(std::net::IpAddr::V4(Ipv4Addr::LOCALHOST), port.server_port ));
+ info!("3");
tokio::spawn(async move {
if let Ok(_) = udp_bind_connect(sid, mux.clone(), &mut session_rx, target_addr ).await {
relay_forwarded_udp(sid, stream, mux_clone, session_rx, port.server_port).await;
diff --git a/src/mux.rs b/src/mux.rs
index 63c16c6..f5c5e84 100644
--- a/src/mux.rs
+++ b/src/mux.rs
@@ -141,7 +141,7 @@ impl MuxHandle {
/// Convenience: send a typed frame.
pub async fn send(&self, frame_type: FrameType, session_id: u32, payload: Vec<u8>) {
// info!("send frame");
- self.send_frame(&Frame::new(frame_type, session_id, payload));
+ self.send_frame(&Frame::new(frame_type, session_id, payload)).await;
}
/// Dispatch an incoming frame to the appropriate session.
diff --git a/src/relay.rs b/src/relay.rs
index 36507bf..634e9ab 100644
--- a/src/relay.rs
+++ b/src/relay.rs
@@ -42,7 +42,7 @@ pub async fn relay_bidirectional_tcp(
match tcp_read.read(&mut buf).await {
Ok(0) => break,
Ok(n) => {
- mux_fwd.send(FrameType::Data, sid, buf[..n].to_vec());
+ mux_fwd.send(FrameType::Data, sid, buf[..n].to_vec()).await;
}
Err(e) => {
debug!("[{}] TCP read error: {}", sid, e);
@@ -77,8 +77,8 @@ pub async fn relay_bidirectional_tcp(
info!("[{}] drop connection", sid);
- mux.send(FrameType::Close, sid, Vec::new());
- mux.drop_session(sid);
+ mux.send(FrameType::Close, sid, Vec::new()).await;
+ mux.drop_session(sid).await;
}
@@ -130,7 +130,7 @@ pub async fn relay_bidirectional_udp_client_side(
let mut a = client_local_port_1.lock().await; *a = Some(addr.port());
- mux_fwd.send(FrameType::Data, sid, buf[..n].to_vec());
+ mux_fwd.send(FrameType::Data, sid, buf[..n].to_vec()).await;
}
Err(e) => {
debug!("[{}] UDP read error: {}", sid, e);
@@ -204,8 +204,8 @@ pub async fn relay_bidirectional_udp_client_side(
_ = break_connection_tcp_check => {}
}
- mux.send(FrameType::Close, sid, Vec::new());
- mux.drop_session(sid);
+ mux.send(FrameType::Close, sid, Vec::new()).await;
+ mux.drop_session(sid).await;
}
pub async fn relay_bidirectional_udp_server_side(
sid: u32,
@@ -236,7 +236,7 @@ pub async fn relay_bidirectional_udp_server_side(
let mut packet = new_udp_header(addr).expect("cannot wrap udp packet");
packet.extend_from_slice(&buf[..n]);
- mux_fwd.send(FrameType::Data, sid, packet.to_vec());
+ mux_fwd.send(FrameType::Data, sid, packet.to_vec()).await;
} else {
warn!("packet came from illegal server location")
}
@@ -301,8 +301,8 @@ pub async fn relay_bidirectional_udp_server_side(
_ = rns_to_udp => {},
}
- mux.send(FrameType::Close, sid, Vec::new());
- mux.drop_session(sid);
+ mux.send(FrameType::Close, sid, Vec::new()).await;
+ mux.drop_session(sid).await;
}
/// udp must still have an associated tcp connection to detect when the connection is over.
@@ -335,7 +335,7 @@ pub async fn relay_forwarded_tcp(
Ok(0) => break,
Ok(n) => {
- mux_fwd.send(FrameType::Data, sid, buf[..n].to_vec());
+ mux_fwd.send(FrameType::Data, sid, buf[..n].to_vec()).await;
}
Err(e) => {
debug!("[{}] TCP read error: {}", sid, e);
@@ -368,8 +368,8 @@ pub async fn relay_forwarded_tcp(
_ = rns_to_tcp => {},
}
- mux.send(FrameType::Close, sid, Vec::new());
- mux.drop_session(sid);
+ mux.send(FrameType::Close, sid, Vec::new()).await;
+ mux.drop_session(sid).await;
}
pub async fn relay_forwarded_udp(
@@ -397,7 +397,7 @@ pub async fn relay_forwarded_udp(
let mut packet = new_udp_header(localhost_server_addr.clone() ).expect("cannot wrap udp packet");
packet.extend_from_slice(&buf[..n]);
- mux_fwd.send(FrameType::Data, sid, packet.to_vec());
+ mux_fwd.send(FrameType::Data, sid, packet.to_vec()).await;
}
Err(e) => {
debug!("[{}] TCP read error: {}", sid, e);
@@ -447,8 +447,8 @@ pub async fn relay_forwarded_udp(
_ = rns_to_tcp => {},
}
- mux.send(FrameType::Close, sid, Vec::new());
- mux.drop_session(sid);
+ mux.send(FrameType::Close, sid, Vec::new()).await;
+ mux.drop_session(sid).await;
}
diff --git a/src/server.rs b/src/server.rs
index 1bf7f10..2bc242b 100644
--- a/src/server.rs
+++ b/src/server.rs
@@ -144,7 +144,7 @@ pub async fn run_server(identity_path: Option<&str>, filter_config: FilterConfig
);
let mux = MuxHandle::new(Arc::clone(&node));
- mux.set_link_id(link_id);
+ mux.set_link_id(link_id).await;
link_muxes.lock().unwrap().insert(link_id, mux);
}
@@ -190,12 +190,12 @@ pub async fn run_server(identity_path: Option<&str>, filter_config: FilterConfig
FrameType::ConnectErr,
sid,
b"invalid payload".to_vec(),
- );
+ ).await;
}
}
FrameType::Data | FrameType::Close => {
// info!("{:?}", frame);
- mux.dispatch(frame);
+ mux.dispatch(frame).await;
}
_ => {}
}
@@ -223,15 +223,15 @@ async fn handle_server_session_tcp(
Ok(s) => s,
Err(e) => {
warn!("[{}] Connection failed: {}", sid, e);
- mux.send(FrameType::ConnectErr, sid, e.to_string().into_bytes());
- mux.drop_session(sid);
+ mux.send(FrameType::ConnectErr, sid, e.to_string().into_bytes()).await;
+ mux.drop_session(sid).await;
return;
}
};
info!("{} how it went {:?}", sid ,stream);
// Signal success
- mux.send(FrameType::ConnectOk, sid, Vec::new());
+ mux.send(FrameType::ConnectOk, sid, Vec::new()).await;
// Data relay (shared implementation)
info!("{} streaming", sid);
@@ -239,8 +239,8 @@ async fn handle_server_session_tcp(
info!("[{}] TCP Closed", sid);
} else {
warn!("[{}] invalid ip address: {:?}", sid, &addr);
- mux.send(FrameType::ConnectErr, sid, "invalid ip address".to_string().into_bytes());
- mux.drop_session(sid);
+ mux.send(FrameType::ConnectErr, sid, "invalid ip address".to_string().into_bytes()).await;
+ mux.drop_session(sid).await;
return;
}
// Attempt TCP connection
@@ -264,8 +264,8 @@ async fn handle_server_session_udp(
Ok(s) => s,
Err(e) => {
warn!("[{}] udp bind failed 1: {}", sid, e);
- mux.send(FrameType::ConnectErr, sid, e.to_string().into_bytes());
- mux.drop_session(sid);
+ mux.send(FrameType::ConnectErr, sid, e.to_string().into_bytes()).await;
+ mux.drop_session(sid).await;
return;
}
};
@@ -273,7 +273,7 @@ async fn handle_server_session_udp(
info!("successfully made connection?");
// Signal success
- mux.send(FrameType::ConnectOk, sid, Vec::new());
+ mux.send(FrameType::ConnectOk, sid, Vec::new()).await;
// Data relay (shared implementation)
relay_bidirectional_udp_server_side(sid, socket, mux, session_rx, filter_config).await;
Served by rngit 1.5.1 - Generated in 0.03s